Skip to content

fix: keep colliding virtual nodes in ConsistentHash - #3611

Merged
He-Pin merged 2 commits into
apache:mainfrom
pjfanning:fix-consistent-hash-collisions-3279
Oct 10, 2026
Merged

He-Pin merged 2 commits into
apache:mainfrom
pjfanning:fix-consistent-hash-collisions-3279

Conversation

@pjfanning

@pjfanning pjfanning commented Oct 9, 2026 •

Copy link
Copy Markdown
Member

Motivation

ConsistentHash stores ring positions in a SortedMap[Int, T]. When virtual nodes of two different nodes hash to the same position, the later entry silently overwrites the earlier one. As a result:

  • the ring depends on the order in which nodes were given or added, so two instances built from the same nodes can route the same key to different nodes
  • removing a node also deletes any position it collides on, even if another node owns that position (or the removed node was never in the ring), so the remaining node loses virtual nodes

Modification

  • Keep every node that claims a ring position: the owner stays in the ring, and the others go in a collisions map, which is normally empty.

  • On a collision, the node with the lowest toString owns the position, independent of insertion order. Nodes are already identified by toString, per the class docs.

  • :- only releases the removed node's claims, and promotes the next claimant if the owner is removed.

  • Documented the collision rule in the class Scaladoc.

  • Added ConsistentHashSpec, using the two node names from the issue whose virtual nodes are known to collide with a factor of 10.

  • apply builds the ring in one sorted batch: each virtual node is encoded as (ring position << 32 | index) in a primitive Long array and sorted once. Single-claimant positions go through a TreeMap builder, and only colliding runs go through the owner rule shared with :+.

Performance

JMH, 10 virtual nodes per node, old = current main, new = this PR:

Operation Nodes Old alloc New alloc Old time New time
nodeFor 10 / 100 / 1000 24 B 24 B 0.08 / 0.17 / 0.21 us 0.08 / 0.11 / 0.21 us
apply 10 24.4 KB 18.5 KB 12 us 22 us (±20, noisy)
apply 100 230 KB 201 KB 219 us 235 us
apply 1000 2.39 MB 1.87 MB 3378 us 2899 us
:+ / :- 10 / 100 / 1000 16 KB / 125 KB / 1.31 MB within ~10% within noise

nodeFor, which is on the per-message path, is unchanged. apply, :+ and :- only run when the set of nodes or routees changes.

Result

The ring, and so nodeFor, is the same regardless of node order, and ch :- node routes like a ring built without that node.

Compared with the current release, routing only changes for keys that land on a colliding ring position. During a rolling upgrade, nodes on old and new versions could route those few keys differently.

Tests

  • actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec without the fix: 5 of 6 failed (the remaining one is an "empty after removing all nodes" sanity check)
  • actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec org.apache.pekko.routing.ConsistentHashingRouterSpec with the fix: 10 passed on 2.13.18 and 3.3.8, including a check that apply builds the same ring as adding the nodes one by one with :+
  • JMH ConsistentHashBench (local only, not committed) comparing main and this PR, results above
  • actor/mimaReportBinaryIssues: passed on 2.13.18 and 3.3.8
  • scalafmt run on changed files

References

Fixes #3279

Motivation:
`ConsistentHash` stores ring positions in a `SortedMap[Int, T]`. When
virtual nodes of two different nodes hash to the same position, the
later entry silently overwrites the earlier one. The ring then depends
on the order in which nodes were given or added, so two instances built
from the same nodes can route the same key to different nodes.
Removing a node also deletes any position it collides on, even if that
position is owned by another node (or the removed node was never in
the ring), so the remaining node loses virtual nodes.

Modification:
- Keep all nodes that claim a ring position: the owner in the ring,
  the others in a `collisions` map, which is normally empty.
- On a collision the node with the lowest `toString` owns the position,
  independent of insertion order.
- `:-` only releases the removed node's claims and promotes the next
  claimant if the owner is removed.
- Document the collision rule in the class Scaladoc.
- Add `ConsistentHashSpec` using two node names whose virtual nodes are
  known to collide.

Result:
The ring, and so `nodeFor`, is the same regardless of node order, and
`ch :- node` routes like a ring built without that node. Routing only
changes compared to before for keys that land on a colliding position.

Tests:
- `actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec`
  without the fix: 5 of 6 failed
- `actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec
  org.apache.pekko.routing.ConsistentHashingRouterSpec` with the fix:
  9 passed on 2.13.18 and 3.3.8
- `actor/mimaReportBinaryIssues`: passed on 2.13.18 and 3.3.8
- scalafmt run on changed files

References:
Fixes apache#3279
Motivation:
Resolving collisions by claiming one virtual node at a time made
`ConsistentHash.apply` insert into the persistent ring map point by
point. Against the previous bulk `SortedMap.empty ++ pairs` build that
is 2-3x the allocation and several times the time (JMH, 10 virtual
nodes per node: 100 nodes 230 KB -> 576 KB, 1000 nodes 2.4 MB -> 7.0 MB).
`apply` is used to rebuild the ring on membership changes by the classic
consistent hashing router, `ConsistentHashingShardAllocationStrategy`
and `ClusterClient`.

Modification:
- Encode each virtual node as `(ring position << 32 | index)` in a
  primitive `Long` array, sort it once, and add positions with a single
  claimant through a `TreeMap` builder.
- Resolve only runs of equal positions (collisions) with the same owner
  rule as `:+`, via a shared `ownerOf` helper.

Result:
`apply` produces the same ring as before this commit, and allocates less
than the original (pre-fix) bulk build at similar speed. JMH, 10 virtual
nodes per node, original vs this commit:
- 10 nodes: 24.4 KB -> 18.5 KB, 12 us -> 22 us (+-20 us, noisy)
- 100 nodes: 230 KB -> 201 KB, 219 us -> 235 us
- 1000 nodes: 2.39 MB -> 1.87 MB, 3378 us -> 2899 us

Tests:
- `actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec
  org.apache.pekko.routing.ConsistentHashingRouterSpec`: 10 passed on
  2.13.18 and 3.3.8, including a new check that `apply` matches adding
  the nodes one by one with `:+`
- `actor/mimaReportBinaryIssues`: passed on 2.13.18 and 3.3.8
- scalafmt run on changed files

References:
Refs apache#3279
@pjfanning

Copy link
Copy Markdown
Member Author

Impact of the routing change for colliding ring positions

Routing only changes for keys that land on a ring position where virtual nodes of two different nodes collide. Everywhere else the ring is identical to main.

How likely a collision is. Ring positions are 32-bit hashes. With P = nodes × virtual nodes positions, the expected number of colliding positions is about P²/2³³ (the default virtual-nodes-factor is 10):

Nodes × virtual nodes Chance of any collision Share of keys that change owner if there is one
10 × 10 ~0.0001% 0.5%
100 × 10 0.01% 0.05%
1000 × 10 1.2% 0.005%
1000 × 100 69% 0.0005%

Only keys on the arc owned by the colliding position (about 1/P of keys) can move, and only when the previous last-written-wins owner differs from the new lowest-toString owner (about half the time).

What this means for each user, including mixed versions during a rolling upgrade:

  • ConsistentHashingShardAllocationStrategy: only the coordinator runs it, and it sorts the nodes, so old and new versions never disagree at the same time. Once the coordinator runs this version, a shard whose id lands on a changed arc may be rebalanced once. For example, with 1000 shards on a 1000 × 10 ring that is about 0.05 shards, and only in the ~1.2% of clusters with a collision.
  • Classic ConsistentHashingRoutingLogic: each sender routes independently, so during an upgrade old and new nodes could send those few keys to different routees, as already happens briefly on any membership change. On main the result depends on routee order, and cluster-aware routers don't sort routees, so two nodes on main can already disagree about colliding keys indefinitely. With this change they agree.
  • ClusterClient receptionist: the ring only chooses which contact points to hand to a client, so a different choice doesn't affect correctness.
  • Typed ConsistentHashingLogic: each router instance is local, so there is no cross-node agreement to break. It uses :-/:+, where main could also drop colliding positions owned by remaining routees when a routee was removed. That is fixed here.

Separately, ch :- node on main can strip virtual nodes from remaining nodes, or change the ring when removing a node that was never in it. Both are fixed here.

Suggested release note: "ConsistentHash now resolves virtual-node hash collisions deterministically and no longer drops virtual nodes; keys at colliding ring positions may map to a different node than before."

@pjfanning pjfanning added this to the 2.0.0-M5 milestone Oct 9, 2026
@pjfanning
pjfanning marked this pull request as ready for review October 9, 2026 16:05
@pjfanning pjfanning added the release notes Need to release note label Oct 9, 2026

@He-Pin He-Pin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Deterministic collision ownership (lowest toString wins) is correct and well covered by the new spec: apply vs incremental equivalence, re-election on removal, and removing a node that is not in the ring. Minor non-blocking nit: the virtualNodesFactor < 1 check added in apply duplicates the one already in the class body.

@He-Pin
He-Pin merged commit eeb3b10 into apache:main Oct 10, 2026
10 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

release notes Need to release note

Projects

None yet

Development

Successfully merging this pull request may close these issues.

ConsistentHash silently drops virtual nodes on hash collisions

2 participants